[spark] Add create_tag_from_watermark procedure - #10198
Open
jackylee-ch wants to merge 2 commits into
Open
jackylee-ch wants to merge 2 commits into
jackylee-ch wants to merge 2 commits into
Conversation
Flink has a create_tag_from_watermark procedure but Spark has none, so a Spark user cannot tag the first snapshot whose watermark reaches a target -- even though Spark already exposes rollback_to_watermark and create_tag_from_timestamp. It is useful when a Flink writer commits event-time watermarks and a Spark job needs a consistent tag at a watermark boundary. Add CreateTagFromWatermarkProcedure, registered as create_tag_from_watermark. It mirrors the existing CreateTagFromTimestampProcedure structure and the Flink procedure's logic: resolve SnapshotManager.laterOrEqualWatermark, fall back to the earliest tag whose watermark already covers the target, else raise SnapshotNotExistException. Tests: CreateTagFromWatermarkProcedureTest commits three snapshots with watermarks 1000/2000/3000 (through the stream write builder, since a Spark batch write carries no watermark) and asserts the resolved snapshot for a below/equal watermark and a SnapshotNotExistException past the last one. Written with Claude Code; verification is mine.
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Purpose
Flink has a
create_tag_from_watermarkprocedure but Spark has none, so aSpark user cannot tag the first snapshot whose watermark reaches a target —
even though Spark already exposes
rollback_to_watermarkandcreate_tag_from_timestamp. It is useful when a Flink writer commitsevent-time watermarks and a Spark job needs a consistent tag at a watermark
boundary.
Change
Add
CreateTagFromWatermarkProcedure, registered ascreate_tag_from_watermark. It mirrors the existingCreateTagFromTimestampProcedurestructure and the Flink procedure's logic:resolve
SnapshotManager.laterOrEqualWatermark, fall back to the earliest tagwhose watermark already covers the target, else raise
SnapshotNotExistException.Tests
CreateTagFromWatermarkProcedureTestcommits three snapshots with watermarks1000/2000/3000 (through the stream write builder, since a Spark batch write
carries no watermark) and asserts the resolved snapshot for a below/equal
watermark and a
SnapshotNotExistExceptionpast the last one.Written with Claude Code; verification is mine.